مستندات فنی CipherMQ
CipherMQ یک بروکر پیام async نوشتهشده به Rust است که پیامهای رمزنگاریشده را میان تولیدکنندهها و مصرفکنندهها، بر پایه هویت گواهی TLS آنها، مسیریابی میکند بدون آنکه محتوای پیام را در هیچ نقطهای رمزگشایی کند.
نمای کلی و معماری
هسته سیستم حول یک مدل آشنا برای هرکسی که با RabbitMQ کار کرده باشد میچرخد:
exchange → binding (routing key) →
queue. اما برخلاف RabbitMQ، CipherMQ از پروتکل استاندارد AMQP استفاده
نمیکند؛ یک پروتکل متنی خطی و اختصاصی روی سوکت TLS خام پیادهسازی شده است (جزئیات در بخش
پروتکل سیمی). تفاوت مهم دیگر این است که بروکر هرگز محتوای پیام را نمیبیند:
پیامها پیش از رسیدن به سرور توسط کلاینت رمزنگاری میشوند و سرور صرفاً یک رلهی رمزنگاریکور
(blind relay) است.
هویت هر کلاینت و بر همین اساس نقش، مجوزها، و سطلهای محدودسازی نرخ او مستقیماً از Common Name (CN) گواهی TLS ارائهشده در زمان دستدهی mTLS استخراج میشود. هیچ لایه جداگانهی نامکاربری/گذرواژه یا توکنی در مسیر داده وجود ندارد.
مدل تحویل
Push-based؛ سرور پیام را فوراً به مصرفکنندگان زنده هل میدهد (نه Pull مانند Kafka)
واحد هویت
Subject CN گواهی X.509 کلاینت (بدون rotation یا revocation در کد بررسیشده)
پایداری صف
فقط در حافظه (VecDeque)؛ فقط متادیتا در PostgreSQL نوشته میشود
همزمانی
DashMap برای صف/binding/consumer + AtomicUsize برای شمارندهها (بدون قفل سراسری)
مدل امنیتی و احراز هویت
۲.۱ احراز هویت متقابل (mTLS)
هر دو listener (بروکر اصلی و کنسول مدیریتی) با rustls::ServerConfig و
WebPkiClientVerifier پیکربندی میشوند تا گواهی کلاینت را در برابر یک
RootCertStore ساختهشده از ca_cert_path اعتبارسنجی کنند.
اتصالی که گواهی معتبر امضاشده توسط همان CA ارائه ندهد، در همان مرحله دستدهی TLS رد میشود —
پیش از آنکه حتی یک بایت دادهی برنامه رد و بدل شود.
پس از دستدهی موفق، CN از فیلد Subject گواهی (OID 2.5.4.3) با
x509-parser استخراج میشود و بهعنوان client_id
در سراسر سیستم استفاده میشود.
کنسول مدیریتی از همان گواهی/کلید/CA سرور اصلی استفاده میکند (در main.rs
مقادیر config.tls.* مستقیماً به admin_server::serve
پاس داده میشوند). به عبارت دیگر، مرز اعتماد TLS برای پورت مدیریتی همان CA است که برای کلاینتهای
عادی استفاده میشود؛ هر گواهی امضاشده توسط آن CA میتواند دستدهی TLS را با پورت مدیریتی کامل کند.
تنها لایهای که واقعاً دسترسی به کنسول را محدود میکند، بررسی نقش در لایه برنامه است
(acl_manager.resolve_role(cn) == Role::Admin)، نه یک مرز TLS جداگانه.
۲.۲ کنترل دسترسی مبتنی بر نقش (ACL)
AclManager هر CN را به یکی از چهار نقش نگاشت میکند:
Sender، Receiver، Admin
یا Unknown (بر اساس عضویت در سه فهرست پیکربندیشده
acl.sender_cns، acl.receiver_cns،
acl.admin_cns). نقش Unknown از همان ابتدا و
پیش از رسیدن به حلقه دستورات رد میشود.
| دستور | Sender | Receiver | Admin |
|---|---|---|---|
declare_queue | ✕ | ✓ | ✓ |
declare_exchange | ✕ | ✓ | ✓ |
bind | ✕ | ✓ | ✓ |
publish / publish_batch | ✓ | ✕ | ✓ |
consume | ✕ | ✓ | ✓ |
ack | ✓ | ✓ | ✓ |
register_public_key | ✕ | ✓ | ✓ |
get_public_key | ✓ | ✕ | ✓ |
resend | ✓ | ✓ | ✓ |
heartbeat | ✕ | ✓ | ✓ |
علاوهبر بررسی سطح دستور، برای declare_queue و bind
یک بررسی سطح منبع (check_resource) نیز اجرا میشود: یک Receiver فقط میتواند
صفی را اعلام یا bind کند که نامش با "{cn}_" شروع شود یا دقیقاً
"{cn}_queue" باشد؛ routing key نیز باید با CN شروع شود. در مقابل،
declare_exchange هیچ بررسی مالکیتی ندارد هر Receiver میتواند هر
نام exchange دلخواهی را اعلام کند.
۲.۳ سرور مدیریتی
admin_server.rs پس از دستدهی TLS، بلافاصله CN را استخراج و نقش را
حل میکند؛ اگر نقش دقیقاً Admin نباشد، اتصال بیصدا بسته میشود
(بدون ارسال پیام خطا به کلاینت). پارامتر پیکربندی admin.allowed_cns
نیز به تابع serve پاس داده میشود اما در امضای تابع با
_allowed_cns نامگذاری شده و در منطق واقعی مجوزدهی استفاده نمیشود
(جزئیات در بخش محدودیتها).
رمزنگاری: در انتقال، سر تا سر، و در سکون
۳.۱ رمزنگاری در انتقال (TLS)
تمام ترافیک چه بروکر اصلی و چه کنسول مدیریتی از طریق TLS 1.2/1.3 با پیادهسازی rustls
عبور میکند و mTLS برای هر دو طرف اجباری است. گواهی سرور و کلید خصوصی آن با
rustls_pemfile بارگذاری میشوند؛ کلید خصوصی باید در قالب PKCS8 باشد
(pkcs8_private_keys) کلیدهای PKCS1/RSA خام (BEGIN RSA PRIVATE KEY)
شناسایی نمیشوند و باید پیش از استفاده تبدیل شوند.
۳.۲ رمزنگاری سر تا سر پیام (End-to-End)
بروکر بهطور طراحیشده محتوای پیام را نمیبیند. ساختار EncryptedInputData
که ناشران ارسال میکنند فقط شامل فیلدهای زیر است:
pub struct EncryptedInputData {
pub message_id: String,
pub receiver_client_id: String,
pub enc_session_key: String, // کلید نشست، رمزشده برای گیرنده
pub nonce: String,
pub ciphertext: String, // بدنه واقعی پیام، هرگز plaintext نیست
}
وجود enc_session_key در کنار nonce و
ciphertext نشاندهنده یک الگوی envelope encryption در سمت کلاینت
است: یک کلید نشست متقارن یکبارمصرف، بدنه پیام را رمز میکند و خودِ آن کلید نشست با کلید عمومی گیرنده
رمز میشود. بروکر نه کلید نشست را میتواند بگشاید و نه بدنه پیام را صرفاً بایتها را ذخیره
(در حافظه) و منتقل میکند.
برای این مدل، بروکر یک دایرکتوری ساده از کلید عمومی ارائه میدهد:
register_public_key (فقط Receiver، همیشه روی client_id
خودِ فراخوان) و get_public_key (فقط Sender، بدون محدودیت روی اینکه
کلید عمومی چه کسی خوانده شود که با توجه به ماهیت «عمومی» بودن این کلید، منطقی است).
۳.۳ رمزنگاری در سکون (At Rest)
کلیدهای عمومی ثبتشده پیش از نوشتن در PostgreSQL با AES-256-GCM (کتابخانه aes_gcm)
رمز میشوند: یک nonce تصادفی ۹۶ بیتی برای هر نوشتن تولید میشود، و ciphertext/tag/nonce بهصورت جداگانه
و base64-encoded در جدول public_keys ذخیره میشوند. کلید رمزنگاری از
encryption.aes_key در فایل پیکربندی (base64، دقیقاً ۳۲ بایت پس از دیکد) خوانده
میشود و برای کل سرور مشترک است نه بهازای هر کلاینت یا هر tenant.
توجه: بدنه پیامها (ciphertext پیام) هرگز در دیتابیس ذخیره نمیشود فقط متادیتای پیام (شناسه، فرستنده، exchange، routing key، زمانها) پایدار میگردد. جزئیات و پیامدهای این موضوع در بخش تضمینهای تحویل شرح داده شده است.
پروتکل سیمی و مرجع دستورات
برخلاف RabbitMQ که از AMQP 0-9-1 استفاده میکند، CipherMQ یک پروتکل متنی خطی اختصاصی روی
سوکت TLS خام پیادهسازی کرده است. هر درخواست یک رشته UTF-8 بهشکل
COMMAND arg1 arg2 ... argN است که با جداکنندهی فاصله تفکیک میشود؛
استخراج آرگومانها معمولاً با splitn انجام میشود، بهاین معنا که فقط
آخرین آرگومان (مثلاً بدنه JSON در publish) میتواند خودش شامل
فاصله باشد نام exchange، صف، یا routing key نباید فاصله داشته باشند.
۴.۱ دستورات بروکر اصلی
| دستور | نقش مجاز | پاسخ موفق | توضیح |
|---|---|---|---|
declare_queue <name> | Receiver, Admin | Queue declared | نیازمند تطابق نام با پیشوند CN |
declare_exchange <name> | Receiver, Admin | Exchange declared | بدون بررسی مالکیت |
bind <queue> <exchange> <key> | Receiver, Admin | Queue bound | تطابق routing key دقیق (بدون wildcard) |
publish <exch> <key> <json> | Sender, Admin | ACK <id> | json باید معادل EncryptedInputData باشد |
publish_batch <exch> <key> <json[]> | Sender, Admin | یک خط به ازای هر پیام | بدون تراکنش اتمیک بین پیامهای دسته |
consume <queue> | Receiver, Admin | (بدون پاسخ فوری) | پیامها async با پیشوند Message: میرسند |
ack <message_id> | Sender, Receiver, Admin | ACK confirmed <id> | بدون بررسی مالکیت پیام |
resend <message_id> | Sender, Receiver, Admin | Resend requested | فقط بررسی مالکیت؛ ارسال مجدد محتوا پیاده نشده |
register_public_key <b64> | Receiver, Admin | Public key registered | همیشه روی client_id خودِ فراخوان |
get_public_key <client_id> | Sender, Admin | Public key: ... | بدون محدودیت خواندن |
heartbeat | Receiver, Admin | Heartbeat received | Sender این دستور را ندارد |
quit / exit | همه | بستن اتصال | خارج از بررسی ACL |
۴.۲ دستورات کنسول مدیریتی
| دستور | خروجی |
|---|---|
json | وضعیت کامل سرور بهصورت یک شیء JSON (شمارندهها، صفها، bindingها، آمار اتصال) |
metrics | متن سبک Prometheus (HELP/TYPE + مقدار برای هر متریک) |
status | خلاصهی خوانا برای انسان |
connections | تعداد اتصال فعال بهتفکیک IP و CN |
queues | تعداد صف/پیام/مصرفکننده |
quit / exit | بستن اتصال |
چرخه حیات پیام
۵.۱ انتشار (Publish)
- بررسی
message_status.contains_key(message_id)؛ در صورت وجود، پیام بهعنوان تکراری رد میشود (idempotency در سطح شناسه پیام). - ساخت رکورد
MessageMetadataباsent_time(UTC RFC3339) و نوشتن همزمان (synchronous) آن در PostgreSQL از طریقstorage.save_metadataاین نوشتن پیش از پاسخ ACK به ناشر تکمیل میشود. - یافتن صفهای مقصد: بررسی
bindings[exchange]و فیلتر بر اساس تطابق دقیقrouting_key(بدون الگوی wildcard مانند#یا*در AMQP). - اگر هیچ صفی matched نشود، خطا بازگردانده میشود با اینحال رکورد متادیتا از قدم قبل همچنان در دیتابیس باقی میماند.
- برای هر صف مقصد: در صورت رسیدن به
max_queue_size، قدیمیترین پیام آن صف حذف (evict) میشود؛ سپس پیام جدید در انتهایVecDequeصف قرار میگیرد. - مصرفکنندگان زنده صف (از طریق کانال
mpsc::UnboundedSender) مطلع میشوند؛ فرستندههای بستهشده در همین گذر پاکسازی میشوند. - ACK نهایی (
ACK <message_id>) به ناشر بازگردانده میشود.
نمودار ۲ مسیر publish. نوشتن متادیتا (نه بدنه پیام) پیش از ACK به ناشر تضمین میشود.
۵.۲ مصرف و تایید (Consume & Ack)
دستور consume صف را (در صورت نبود) اعلام میکند، فرستنده کانال کلاینت را
به فهرست مصرفکنندگان آن صف اضافه میکند، و بدون ارسال هیچ پاسخ فوری به کلاینت پس از ۱۰۰
میلیثانیه تاخیر (برای پرهیز از race با پاسخ ثبت مصرفکننده)، پیامهای در انتظارِ قابلتحویل صف را
با فاصله ۱۰ میلیثانیه بین هر پیام برای مصرفکننده جدید بازپخش میکند.
نمودار ۳ مسیر consume/ack. توجه شود که ack، پیام را از تمام صفهای مقصدش حذف میکند، نه فقط از صفی که این مصرفکننده روی آن ثبت شده.
تضمینهای تحویل پیام
با جمعبندی رفتارهای بخشهای قبل، تضمین تحویل CipherMQ را میتوان اینگونه خلاصه کرد:
| جنبه | تضمین واقعی |
|---|---|
| تحویل حین اجرا (process زنده) | تقریباً «حداقل یکبار» (at-least-once) تا زمانی که صف سرریز نشود و پیام acknowledge نشده باشد |
| پایداری در برابر ریاستارت بروکر | ❌ فقط متادیتا در PostgreSQL میماند؛ بدنه رمزشده پیام (که در حافظه است) از بین میرود |
| idempotency انتشار مجدد | بر اساس message_id، تا زمانی که رکورد message_status آن به دلیل ack یا سرریز صف پاک نشده باشد |
| استقلال fan-out بین چند صف | ❌ یک ack یا یک سرریز، روی همهی صفهای مقصد اثر سراسری دارد |
| ترتیب تحویل در یک صف | FIFO (صفبندی با VecDeque) |
| تلاش مجدد تحویل به سوکت | حداکثر ۳ تلاش با backoff ۱۰ میلیثانیهای، سپس قطع اتصال مصرفکننده |
چون بدنه پیام هرگز در دیتابیس نوشته نمیشود، دستور resend که در
هنگام راهاندازی سرور برای پیامهای تاییدنشده لاگ میشود نمیتواند خودِ محتوای پیام را از
سمت بروکر بازیابی کند؛ طراحی موجود صرفاً یک بررسی مالکیت و یک اعلام «درخواست ارسال مجدد پذیرفته
شد» است. بازارسال واقعی محتوا باید توسط ناشر اصلی (که هنوز نسخهی رمزشده پیام را دارد) انجام شود.
مدیریت وضعیت و همزمانی
ServerState بهطور کامل بدون یک قفل سراسری واحد ساخته شده است؛ هر
نگاشت (صفها، bindingها، exchangeها، مصرفکنندگان، وضعیت پیامها) یک DashMap
مجزاست که قفل را در سطح shard نگه میدارد، و شمارنده اتصال فعال یک AtomicUsize
است نه RwLock<usize>.
cleanup_old_messages
هر ۶۰ ثانیه اجرا میشود؛ پیامهای تحویلدادهشده اما تاییدنشدهای که از delivered_time آنها بیش از max_age (پیشفرض ۳۰۰۰ ثانیه ≈ ۵۰ دقیقه) گذشته باشد را از صف و نگاشتهای ردیابی حذف میکند.
cleanup_dead_consumers
هر ۶۰ ثانیه، فرستندههای mpsc بستهشده را از فهرست مصرفکنندگان هر صف پاک میکند و صفهای مصرفکنندهی خالیشده را حذف میکند.
HeartbeatMonitor
هر ۱۰ ثانیه بررسی میکند و کلاینتهایی که در بازه timeout_duration پیامی نفرستادهاند را از نگاشت heartbeat حذف میکند (بدون بستن فعال اتصال TCP آنها).
Storage flush loop
هر ۵۰۰ میلیثانیه یا هنگام رسیدن به ۵۰ آیتم بافرشده، بهروزرسانیهای delivered/acknowledged را بهصورت دستهای در PostgreSQL مینویسد.
نمودار ۴ وضعیتهای ممکن یک پیام از دید message_status.
لایه ذخیرهسازی
Storage از یک الگوی actor استفاده میکند: یک تسک پسزمینهی واحد،
مالکیت انحصاری PgPool را در اختیار دارد و تمام عملیات دیتابیس از طریق
کانال mpsc::channel::<StorageCommand> به آن ارسال میشوند نه از
طریق دسترسی مستقیم و همزمان به pool از چندین تسک.
- Pool: حداکثر ۲۰۰ / حداقل ۱۰ اتصال، طول عمر حداکثر ۱ ساعت، idle timeout ۱۰۰ دقیقه، acquire timeout ۵۰ دقیقه.
- اتصال اولیه: تا ۵ تلاش با فاصله ۵ ثانیهای؛ در صورت شکست همه، راهاندازی سرور با خطا متوقف میشود.
- بازیابی حین اجرا: اگر
pool.is_closed()شود، تسک storage پیش از پردازش دستور بعدی، یک تلاش اتصال مجدد انجام میدهد. - دستهبندی نوشتن (batching): بهروزرسانیهای
delivered_timeوacknowledged_timeدر یکHashMapبافر میشوند و هر ۵۰۰ میلیثانیه یا در رسیدن به ۵۰ آیتم، با یک کوئری دستهای (UPDATE ... FROM unnest($1::text[])) نوشته میشوند برخلاف نوشتن اولیه متادیتا (SaveMetadata) که همزمان و بلافاصله انجام میشود، چون ACK به ناشر به آن وابسته است.
| جدول | ستونها | هدف |
|---|---|---|
message_metadata | message_id (PK), client_id, exchange_name, routing_key, sent_time, delivered_time, acknowledged_time | ردیابی و audit trail بدون بدنه پیام |
public_keys | client_id (PK), public_key_ciphertext, nonce, tag | دایرکتوری کلید عمومی، رمزشده با AES-256-GCM سرور |
محدودسازی نرخ و محافظت در برابر اضافهبار
الگوریتم پایه token bucket کلاسیک است: در هر take()، تعداد توکنها
بر اساس زمان سپریشده از آخرین شارژ و نرخ پیکربندیشده (rps) تا سقف
burst بازپر میشود؛ اگر حداقل یک توکن موجود باشد، مصرف و درخواست پذیرفته
میشود.
سطل global
بهازای هر CN، برای همه دستورات بررسی میشود (global_rps / global_burst).
سطل publish
سطلی جداگانه و سختگیرانهتر، فقط برای publish/publish_batch، علاوهبر سطل global (هر دو باید عبور کنند).
محدودیت اتصال
تعداد اتصال همزمان بهازای IP و بهازای CN، در لحظه پذیرش اتصال (پیش از ورود به حلقه دستورات) بررسی میشود.
تابآوری و مدیریت خطا
خاموشی آرام (Graceful Shutdown)
دریافت SIGTERM/SIGINT (یا Ctrl+C در غیر یونیکس) از طریق tokio::select! و الگوی Notify؛ حلقه پذیرش اتصال جدید متوقف و پاکسازی نهایی اجرا میشود.
timeout خواندن سوکت
هر خواندن از کلاینت با tokio::time::timeout و سقف ۶۰۰ ثانیه محافظت میشود؛ عبور از این زمان باعث قطع نمیشود، فقط چرخه select! را برای بار دیگر تکرار میکند.
تلاش مجدد تحویل
نوشتن پیام روی سوکت مصرفکننده تا ۳ بار با backoff ۱۰ میلیثانیهای تکرار میشود؛ شکست نهایی باعث قطع اتصال آن مصرفکننده میشود (پیام همچنان در صف باقی میماند تا مصرفکننده بعدی).
اتصال مجدد دیتابیس
در راهاندازی: ۵ تلاش با فاصله ۵ ثانیه. حین اجرا: تشخیص pool بستهشده و یک تلاش اتصال مجدد پیش از پردازش دستور بعدی.
همچنین، در main.rs، هنگام شکست دستدهی TLS، state.decrement_client_count()
فراخوانی میشود درحالیکه increment_client_count() فقط پس از عبور
موفق از بررسی نقش و rate limiter در handle_client اجرا میشود. این
عدم تقارن بهمعنای کاهش شمارنده بدون افزایش متناظر است.
سرور مدیریتی، متریکها و لاگ
کنسول مدیریتی (پیشفرض 127.0.0.1:9091، غیرفعال مگر
admin.enabled = true) یک پروتکل متنی مشابه بروکر اصلی ولی جداگانه
ارائه میدهد و صرفاً برای نقش Admin در دسترس است.
۱۱.۱ متریکها
خروجی metrics در قالب سبک Prometheus (# HELP / # TYPE) است:
| متریک | نوع |
|---|---|
ciphermq_messages_published_total | counter |
ciphermq_messages_acked_total | counter |
ciphermq_messages_delivered_total | counter |
ciphermq_errors_total | counter |
ciphermq_connections_total | counter |
ciphermq_connections_rejected_rate_limit_total | counter |
ciphermq_connections_rejected_conn_limit_total | counter |
ciphermq_requests_rate_limited_total | counter |
ciphermq_connected_clients / queue_count / messages_in_queues / consumer_count | gauge (از snapshot دورهای) |
ciphermq_uptime_seconds | gauge |
دستور json یک تصویر کاملتر ارائه میدهد که شامل هم شمارندههای
اتمیک زنده (live) و هم آخرین ServerSnapshot
ذخیرهشده است این دو ممکن است به دلیل ماهیت بهروزرسانی دورهای snapshot، لحظهای با هم یکسان نباشند.
۱۱.۲ لاگگیری
بر پایه tracing + tracing-subscriber: سه
فایل جداگانه و چرخشی (rolling::hourly/daily/never بر اساس
logging.rotation) برای سطوح info/debug/error بهصورت JSON، بهعلاوه یک
لایهی «pretty» روی stdout که با RUST_LOG یا logging.level
فیلتر میشود.
پیکربندی (config.toml)
پیکربندی با toml + serde بارگذاری و اعتبارسنجی میشود؛ نمونهای معادل بخشهای تعریفشده در config.rs:
# اتصال اصلی بروکر [server] address = "0.0.0.0:5671" connection_type = "tls" # فعلاً تنها مقدار پشتیبانیشده [tls] cert_path = "/etc/ciphermq/certs/server.crt" key_path = "/etc/ciphermq/certs/server.key" # باید PKCS8 باشد ca_cert_path = "/etc/ciphermq/certs/ca.crt" [logging] level = "info" rotation = "daily" # hourly | daily | هر مقدار دیگر = never info_file_path = "logs/info.log" debug_file_path = "logs/debug.log" error_file_path = "logs/error.log" max_size_mb = 50 [database] host = "127.0.0.1" port = 5432 user = "ciphermq" password = "********" dbname = "ciphermq" [encryption] algorithm = "x25519_chacha20_poly1305" # فقط اعتبارسنجی میشود، رجوع به بخش رمزنگاری aes_key = "<base64 دقیقاً معادل ۳۲ بایت>" [performance] max_queue_size = 10000 heartbeat_timeout = 60 cleanup_interval = 3000 consumer_cleanup_interval = 60 idle_timeout_sec = 60 [admin] enabled = true address = "127.0.0.1:9091" allowed_cns = [] # بارگذاری میشود اما در مجوزدهی واقعی استفاده نمیشود [rate_limit] enabled = true global_rps = 100.0 global_burst = 200.0 publish_rps = 50.0 publish_burst = 100.0 max_connections_per_ip = 50 max_connections_per_cn = 20 [acl] sender_cns = ["publisher-svc-1"] receiver_cns = ["worker-svc-1"] admin_cns = ["ops-admin"]
قوانین اعتبارسنجی اجباری در Config::load:
- اگر
connection_type = "tls"، هر سه مسیرtls.*الزامیاند. database.host،database.user،database.dbnameنباید خالی باشند.encryption.algorithmباید دقیقاً"x25519_chacha20_poly1305"باشد.encryption.aes_keyباید base64 معتبر و پس از دیکد دقیقاً ۳۲ بایت باشد.